Repository navigation
feat(connectors): add sink_template and source_template starting points - #3965
ryerraguntla wants to merge 185 commits into
Conversation
There was a problem hiding this comment.
Pull request overview
Warning
Copilot couldn't run its full agentic review because it didn't start before the timeout. Make sure your repository has a runner available, or add a copilot-code-review.yml file specifying one with the runs-on attribute. See the docs for more details.
Adds ready-to-copy “template” sink and source connector crates (plus docs and skill checklists) to standardize connector implementations and reduce recurring review feedback (secrets handling, retries/circuit-breaker, config validation, cursor staging, and canonical tests).
Changes:
- Introduces
sink_templateandsource_templatecrates implementing the expected connector scaffolding withTODO(ConnectorDeveloper)placeholders for backend-specific logic. - Adds example runtime configs and README/docs entries to make the templates discoverable and easy to adopt.
- Updates
.claude/skillsguidance (including newconnector-pr-review) to encode recurring review blockers and conventions.
Reviewed changes
Copilot reviewed 21 out of 22 changed files in this pull request and generated 7 comments.
Show a summary per file
| File | Description |
|---|---|
| core/connectors/sources/source_template/src/lib.rs | New source template implementation + cursor staging + tests |
| core/connectors/sources/source_template/config.toml | Example source template connector config |
| core/connectors/sources/source_template/README.md | Source template usage guidance and expectations |
| core/connectors/sources/source_template/Cargo.toml | Source template crate metadata/deps for workspace + cdylib |
| core/connectors/sources/README.md | Registers source_template in sources index |
| core/connectors/sinks/sink_template/src/lib.rs | New sink template implementation + batching/circuit breaker + tests |
| core/connectors/sinks/sink_template/config.toml | Example sink template connector config |
| core/connectors/sinks/sink_template/README.md | Sink template usage guidance and expectations |
| core/connectors/sinks/sink_template/Cargo.toml | Sink template crate metadata/deps for workspace + cdylib |
| core/connectors/sinks/README.md | Registers sink_template in sinks index |
| core/connectors/runtime/example_config/connectors/source_template.toml | Runtime example config for source template |
| core/connectors/runtime/example_config/connectors/sink_template.toml | Runtime example config for sink template |
| core/connectors/BLOG_POST.md | Draft blog post announcing the templates |
| Cargo.toml | Adds both template crates to workspace members |
| .claude/skills/connectors-overview/SKILL.md | Updates connector overview skill (incl. review checklist pointers) |
| .claude/skills/connector-source/TEMPLATE.md | Reworks source template kit content to align with ACK/NACK + patterns |
| .claude/skills/connector-source/SKILL.md | Documents ACK/NACK contract and staging rules more explicitly |
| .claude/skills/connector-sink/TEMPLATE.md | Reworks sink template kit content and pre-flight guidance |
| .claude/skills/connector-sink/SKILL.md | Adds pointers to PR pre-flight skill |
| .claude/skills/connector-runtime/SKILL.md | Updates metrics naming guidance |
| .claude/skills/connector-pr-review/SKILL.md | New PR review checklist skill for connector PRs |
💡 Add a code-review agent skill or configure MCP servers for context-aware, tailored reviews. Learn more in the docs.
| self.config.retry_max_delay.as_deref(), | ||
| DEFAULT_RETRY_MAX_DELAY, |
| // and unknown keys are rejected outright so a typo in a TOML file fails at | ||
| // load time instead of silently doing nothing. | ||
|
|
||
| #[derive(Debug, Clone, Serialize, Deserialize)] |
| #[serde(serialize_with = "iggy_common::serde_secret::serialize_secret")] | ||
| pub connection_string: SecretString, |
| #[serde( | ||
| default, | ||
| serialize_with = "iggy_common::serde_secret::serialize_optional_secret" | ||
| )] | ||
| pub auth_token: Option<SecretString>, |
| #[tokio::test] | ||
| async fn poll_returns_empty_without_error_when_circuit_is_open() { | ||
| let source = TemplateSource::new(1, test_config(), None); | ||
| source.circuit_breaker.record_failure().await; | ||
| source.circuit_breaker.record_failure().await; | ||
| source.circuit_breaker.record_failure().await; | ||
| assert!(source.circuit_breaker.is_open().await); | ||
|
|
||
| let result = source | ||
| .poll() | ||
| .await |
Codecov Report✅ All modified and coverable lines are covered by tests. Additional details and impacted files@@ Coverage Diff @@
## master #3965 +/- ##
============================================
+ Coverage 84.76% 84.97% +0.21%
- Complexity 1405 1575 +170
============================================
Files 1224 1236 +12
Lines 177971 180844 +2873
Branches 144285 144501 +216
============================================
+ Hits 150854 153671 +2817
+ Misses 23100 23097 -3
- Partials 4017 4076 +59
🚀 New features to boost your workflow:
|
Co-authored-by: Copilot Autofix powered by AI <175728472+Copilot@users.noreply.github.com>
….com/ryerraguntla/iggy into docs/connector-template_code_blog_post
…apache#3970) Every journal flush issued two serialized fdatasyncs, log then index, so an fsync-gated topic paid two device round trips per ack. Both files now persist under one futures::future::join: 13.4k to 26.1k msg/s and p50 5.83 to 3.01 ms on ext4 with enforce_fsync and messages_required_to_save = 1. Cursors advance only after both saves succeed, so a failed half keeps its slot and a later flush rewrites the same positions instead of appending a duplicate. Removing the write barrier forced a new recovery contract. An index entry is derived from the log and, without enforce_fsync, proves nothing about it, so recovery checksum-walks every segment from byte 0 and keeps a clean index untouched. Every entry must match its log batch's offset, timestamp, partition and extent, or the index is rebuilt: ascending columns alone let a corrupt entry point at the wrong valid batch and polls silently skipped offsets. Under enforce_fsync serialized flushes prove the prefix below the last entry, so a healthy boot still anchors there; a gap deeper than the one in-flight chunk refuses as FsyncedLogLoss, and a rebuild stopping below the last entry's position refuses as FsyncedRebuildShortfall rather than truncating bytes a completed flush made durable. Probe-budget exhaustion always refuses, since recovering as empty would serve the segment and re-mint offsets over bytes never scanned. Previously any index the log outran refused boot outright and tombstoned single-replica partitions with healthy logs. A persist failure on a cluster-committed op panicked on the shard pump, where compio::runtime::spawn swallows the unwind: the process stayed up with every partition on that shard dead. The partition now fences itself, the pump breaks, flips the shared shutdown flag before its final flush (the flush writes to the failed device and can stall forever; the flag is what arms the watchdog and the bounded drain), and a failed shutdown flush fences its partition too, so the exit is ShardFatal and non-zero rather than a clean report over unpersisted committed data. Also: the active segment fills the same read-state slot sealed segments use, ending one openat per head poll; latency extremes backfill only as a whole group; segment writer save methods are crate-private; index scans yield on a stride.
Merging adjacent producer batches released their byte permits before the write completed, allowing buffered and in-flight data to exceed max_buffer_size. Large permit totals could also overflow during the merge and permanently reduce available capacity. Retain each permit separately until send_internal completes. Validate batch sizes before u32 conversion, split merged writes at the wire limit, and handle zero buffer and in-flight limits without invalid semaphores. Update max_buffer_size documentation and add regression coverage for permit lifetime, unlimited limits, and budgets above u32::MAX. Keep the edge.6 package versions synchronized. Fixes apache#3932
…he#3958) Our partition operations did not use the `request` counter which meant that partition ops did not participate in the deduplication process. This PR makes the partition operations bump the request counter for all of the SDKs
…up (apache#4007) A partition group with no history of its own seeds its view and log_view from the metadata view this replica happens to publish. Nothing bounds that view by the group's real one, so a replica that materialises the group late, after a further metadata election, starts Normal on an empty log above every peer. That empty log is canonical in the next DVC merge: as a backup it collects a nack quorum against committed ops and the group deadlocks for good, the superblock persisting the view so a restart re-wedges; as the primary-by-index it answers the real primary's heartbeat with a StartView at op 0 that resets the peers' heads under their commit floor, and every new write is fenced. Introduced with the seed in apache#3985. Record the view the admitting primary mints the create in on the committed metadata entry and seed from that instead. Every replica reads the same value, and the group's view never drops below it, so an empty log at the creation view is the floor a view-0 start used to be, never canonical over committed history. The reconciler no longer needs the published-view sentinel, its deferral, or the double atomic read that came with it. The value rides the request body, not the prepare header: a header's view is the one field a post-view-change retransmit restamps (replicate_preflight refuses the frame otherwise), so a header-derived value is a per-replica property of delivery history and durably diverges the metadata STM. The body is sealed by checksum_body and never restamped; entries journaled before the trailing field decode it as 0, the safe floor. Chosen over electing into the seed: no extra round trip per group, no consensus change, and a late materialiser joins as a backup at or below its peers instead of forcing an election it may win with an empty log. The snapshot format moves to version 5 with the new trailing field defaulted, so version 3 and 4 checkpoints still decode.
|
moving this to draft to not trigger stale bot |
…#4224) The link shard that owns a peer's replica socket burns a full core under sustained replicated writes. Nothing batches its reads: read_message issues one io_uring read for the 256-byte header and a second for the body. Bytes already sitting in the kernel receive queue wait for the next call. A read-ahead buffer sized by message_bus.replica_read_buffer_size lets one socket read serve a whole burst. It adds no latency, because the reader never waits for the buffer to fill. One fill delivers at most one buffer, so a body above that size crosses the buffer in buffer-sized pieces and pays one extra pass over the payload and one read per piece. A read that large goes straight into the frame's own buffer instead. Zero keeps the unbuffered path, so the A/B baseline is a config change and not a separate build. The plaintext writer already coalesces into one writev, so it is left alone. replica_socket_reads_total and replica_inbound_frames_total carry the evidence on the real workload: their ratio is the batching factor, and only link shards bump them. A read counts once it completes, so a link torn down mid-read cannot inflate the ratio.
Co-authored-by: Copilot Autofix powered by AI <175728472+Copilot@users.noreply.github.com>
# Conflicts: # .claude/skills/connector-source/TEMPLATE.md
apache#4104 landed on master after this branch's commits were written and replaced the connectivity retry API: ConnectivityConfig is gone, check_connectivity_with_retry now takes an owned RetryPolicy. Map the old fields onto the new struct with identical semantics (max_open_retries -> max_attempts, retry_delay -> base_delay, open_retry_max_delay -> max_delay) in both templates. Co-Authored-By: Claude Sonnet 5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01JpGKmaGSG6RdBUiaq3BW2f
….com/ryerraguntla/iggy into docs/connector-template_code_blog_post
|
@ryerraguntla looks like merge is broken in this PR |
|
history broken in this PR, please reopen a new one as fresh |

Fixes #3956
What
Adds two ready-to-copy connector crates under
core/connectors/:core/connectors/sinks/sink_template/core/connectors/sources/source_template/Each implements every framework-level requirement already expected of a production connector (config parsing with unknown-field rejection,
SecretString-wrapped credentials, connectivity validation inopen(), retry/backoff + circuit breaker viaiggy_connector_sdk::retry, batch processing / Ack-Nack cursor staging, and a canonical test suite), leaving only the backend-specific bits markedTODO(ConnectorDeveloper):Also included:
core/connectors/docs/authoring-sinks-and-sources.md— a written guide walking through the required patterns and why each one exists, meant to be cross-linked from the connector READMEs and reused as the basis for an Apache Iggy blog post..claude/skills/connector-*skill definitions, including a newconnector-pr-reviewskill intended to catch the same recurring review issues (missing retry/circuit-breaker, plaintext secrets, unvalidated config, missing cursor staging) automatically on future connector PRs.Cargo.toml,Cargo.lock) and README entries (core/connectors/sinks/README.md,core/connectors/sources/README.md,core/connectors/README.md) for both new crates.core/connectors/runtime/example_config/connectors/{sink,source}_template.toml).Why
Per the discussion in #3956: connector PRs have been hitting the same review comments repeatedly (improperly typed secrets, missing state-staging, inconsistent error handling, thin test coverage). Rather than repeating that feedback PR after PR, these templates make the required pattern the path of least resistance — a new connector starts from working, reviewed code and only needs the backend-specific logic filled in.
Local Execution
Passed
Pre-commit hooks ran
Testing
cargo fmt --allcargo clippy --all-targets --all-features -- -D warningscargo build -p iggy_connector_template_sink -p iggy_connector_template_sourcecargo test -p iggy_connector_template_sink -p iggy_connector_template_sourcecargo machetecargo sort --workspacetypos./scripts/ci/license-headers.sh --checkAI Usage
Claude
Implementation was AI-generated with human review
Each line of the code is reviewed and can be explained by developer